场景案例:如何用 Flink SQL CDC 实现实时数据同步?
0. 引言
业务增长后,一份数据往往需要写入多种数据库:MySQL(主库)、Redis(缓存)、Elasticsearch(查询分析)。如何低成本地完成多库同步?本文从"改业务代码"与"中间件方案"的痛点出发,介绍 Flink SQL CDC 的原理与实现——只需几行 SQL 即可完成全量 + 增量实时同步。
1. 业务场景:多库同步的三种方案
1.1 方案一:业务代码直写多库(不推荐)
在业务代码中同时写入 MySQL、Redis、Elasticsearch。缺点:
- 代码严重耦合:每次修改都要重新测试、上线;
- 性能受损:业务系统多次写不同数据库;
- 全量同步靠人工:增量前需人工做一次全量。
1.2 方案二:消息中间件解耦
优点:业务与写入逻辑解耦、业务只需写一次 Kafka。缺点:仍需埋点改代码、需运维 Kafka、多库多表时开发量激增、全量同步仍需人工。
1.3 方案三:Flink CDC(推荐)
CDC(Change Data Capture,变化数据捕获) 是捕获数据变化的技术——MySQL 主从同步跟随 binlog 就是典型 CDC 场景。Flink CDC 用 Flink 实现 CDC:先全量同步,再跟随 binlog 实时增量同步,且全量/增量封装一体、支持 SQL,无需引入 Kafka、无需改业务代码。
2. 实现原理:快照 + binlog 增量
以 MySQL → 目标库为例,Flink CDC 分两步:全量同步(快照)→ 增量同步(binlog)。
2.1 全量同步:快照流程
关键点:可重复读事务 + 全程只读 → 无幻读 → 扫描数据与 binlog 位置时间点对齐,保证"一致快照"(不多不少、与源库完全相同)。
2.2 增量同步:binlog 跟随 + Checkpoint
全量完成后,跟随源库 binlog 将每次变更实时同步到目标库。Flink 周期性执行 Checkpoint 记录已同步的 binlog 位置,作业故障重启后从最近 Checkpoint 恢复。配合目标写入幂等,整个同步过程可达 exactly once 可靠性。
3. 具体实现
3.1 方式一:基于 DataStream
SourceFunction<String> sourceFunction = MySQLSource.<String>builder()
.hostname("127.0.0.1").port(3306)
.databaseList("db001")
.username("root").password("123456") // 测试用,生产禁用 root 与弱口令
.deserializer(new StringDebeziumDeserializationSchema())
.build();
// Elasticsearch 目标库 Sink(省略构造细节)
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.addSource(sourceFunction) // 源:MySQL(封装全量+增量)
.addSink(esSinkBuilder.build()) // 目标:Elasticsearch
.setParallelism(1); // 保持消息顺序
env.execute("FlinkCdcDemo");MySQLSource 连接器封装了全部复杂操作。缺点:写入 ES 的是整串字符串(未解析字段,查询效率低)、仍需写代码。
3.2 方式二:Flink SQL CDC(三行 SQL)
-- 1. 配置源:MySQL
CREATE TABLE sourceTable (
id INT, name STRING, counts INT, description STRING
) WITH (
'connector' = 'mysql-cdc',
'hostname' = '192.168.1.7', 'port' = '3306',
'username' = 'root', 'password' = '123456',
'database-name' = 'db001', 'table-name' = 'table001'
);
-- 2. 配置目标:Elasticsearch
CREATE TABLE sinkTable (
id INT, name STRING, counts INT
) WITH (
'connector' = 'elasticsearch-7',
'hosts' = 'http://192.168.1.7:9200', 'index' = 'table001'
);
-- 3. 启动同步作业(可指定字段、支持 GROUP BY/Window 等复杂逻辑)
INSERT INTO sinkTable SELECT id, name, counts FROM sourceTable;三个 SQL 完成全量 + 增量实时同步,且目标库字段与 SELECT 指定完全一致。Flink SQL 底层也会转化为 DataStream 执行。
4. 小结
- 三种方案对比:业务直写(耦合高)→ Kafka 中间件(解耦但复杂)→ Flink CDC(无埋点、无 Kafka、全量+增量一体);
- 原理:全局读锁 + 可重复读事务的一致快照 → binlog 位置 → 全表扫描 → Checkpoint 增量续传 → 幂等写入 = exactly once;
- 两种实现:DataStream(灵活,适合 SQL 暂不支持的功能)与 Table & SQL(简单,适合常规同步);
- 生产建议:禁 root 账号、弱口令;DataStream 与 SQL 两种方式都需掌握。
下一章讲解:竟然还有分布式的 JVM?